Skip to content

ZeroMQ 代理心跳检测:让 proxy 感知 server 存活 ​

用 zmq_proxy() 搭建异步客户端/服务器模式的代理时,有一个容易被忽视的坑:当 server 全部退出或网络中断后,管道里积压的未消费消息会让 proxy 产生高额 CPU 开销。解法是让 proxy 能感知 server 是否存活——server 数量归零时及时关闭代理。本文介绍用 ZMTP 心跳与 socket monitor 事件监听组合实现这一检测,附完整可运行的 C++ 示例。适合正在用 libzmq 搭建 ROUTER/DEALER 代理的开发者。

异步客户端/服务器模式与 zmq_proxy ​

ZeroMQ 官方指南的异步客户端/服务器模式用代理连接两侧:多个 client 与多个 server 不直接通信,中间由一个 proxy 转发。

代理本身由 zmq_proxy() 实现:

cpp
int zmq_proxy(const void *frontend, const void *backend, const void *capture);

几个关键行为(依据 libzmq 官方文档):

  • zmq_proxy() 在当前线程中运行,仅在当前上下文关闭时返回(此时返回 -1,errno 为 ETERM 等);
  • 代理将前端连接到后端,数据从前端流向后端只是概念上的——代理是完全对称的,前后端之间没有技术差异;
  • 如果 capture 参数不为 NULL,代理会把前后端收到的所有消息再发送到 capture 套接字,capture 应为 ZMQ_PUB、ZMQ_DEALER、ZMQ_PUSH 或 ZMQ_PAIR 类型,常用于旁路观察流量;
  • 调用前需要先设置必要的套接字选项,并完成 bind 或 connect。

为什么需要检测 server 存活 ​

zmq_proxy() 在管道中存在积压的未消费消息时,会持续空转并产生高额 CPU 开销。导致积压的典型场景有两个:

  1. server 进程退出(正常退出或崩溃);
  2. 网络中断,server 无法触达但连接尚未被判定失效。

这两种场景需要不同的检测手段,区别在于 TCP 连接的回收机制:

  • 进程退出时:TCP 连接信息由内核维护。进程退出后,内核回收其所有 TCP 连接资源,主动发送第一次挥手的 FIN 报文,后续挥手过程也在内核完成,不需要进程参与。也就是说,即使 server 进程已经不存在,两端仍能完成 TCP 四次挥手。此时只需用 socket monitor 监听 ZMQ_EVENT_DISCONNECTED 事件,即可感知断开。
  • 网络中断时:TCP 连接仍然存在,对端崩溃、拔线、断网都不会触发 ZMQ_EVENT_DISCONNECTED。这时需要在应用层发送心跳包,靠超时判定 server 是否存活。

ZMTP 心跳:三个套接字选项 ​

zmq_setsockopt() 中与心跳相关的选项有三个:

选项 语义 依赖与生效条件
ZMQ_HEARTBEAT_IVL 心跳发送间隔(毫秒)。大于 0 时,每隔 IVL 毫秒发送一次 PING ZMTP 命令 独立生效,默认 0(不发心跳)
ZMQ_HEARTBEAT_TIMEOUT 发送 PING 后未收到任何流量时,判定连接超时的等待时长(毫秒) 仅在 IVL 已设置且大于 0 时有效;默认取 IVL 值
ZMQ_HEARTBEAT_TTL 告知远端的心跳超时(毫秒)。远端在 TTL 内未收到任何流量则断开连接 依赖 IVL;内部向下取整到 0.1 秒,小于 100 的值无效,最大 6553599

两个容易误读的细节(均为官方文档明确说明):

  • 取消 TIMEOUT 超时不要求收到 PONG——任何收到的流量都会取消本次超时计时;
  • TTL 是设置在远端生效的:本端通过该选项告知对端「如果 TTL 时间内没有收到任何流量,就由你那边超时断开」。

在代理的后端套接字上添加心跳(需要在 bind 之前设置,套接字选项只影响之后建立的连接):

cpp
_backend.setsockopt(ZMQ_HEARTBEAT_IVL, 1000);
_backend.setsockopt(ZMQ_HEARTBEAT_TIMEOUT, 5000);

_frontend.bind("tcp://*:5570");
_backend.bind("tcp://*:5571");
_capture.bind("inproc://capture");

这样 proxy 会对连接到 backend 的每一个 TCP 连接发送心跳。设置 socket monitor 之后,心跳超时导致的断开同样会触发 ZMQ_EVENT_DISCONNECTED——实测在局域网内两台机器分别启动 proxy 与 server,断开网络连接后,达到超时时间即触发该事件。用 tcpdump -i lo -vnnX port 5571 抓包,可以看到 proxy 与 server 之间一秒一次的 PING/PONG 心跳。

用 socket monitor 捕获断开事件 ​

cpp
int zmq_socket_monitor(void *socket, char *endpoint, int events);

zmq_socket_monitor() 允许应用跟踪某个套接字上的连接事件。每次调用会创建一个 ZMQ_PAIR 套接字并绑定到指定的 inproc:// 端点(监控端点只支持 inproc 传输;被监控的套接字本身须使用面向连接的传输,如 tcp/ipc)。应用自己创建一个 ZMQ_PAIR 套接字连接到该端点,即可收到事件。events 参数是事件位掩码,ZMQ_EVENT_ALL 表示全部事件。

每个事件以两帧发送:

  • 第一帧:事件编号(16 位)+ 事件值(32 位,按事件类型提供附加数据,例如 ZMQ_EVENT_DISCONNECTED 的值是底层网络套接字的文件描述符);
  • 第二帧:受影响端点的字符串。

解析两帧的官方示例代码:

cpp
static int get_monitor_event(void *monitor, int *value, char **address)
{
    // 第一帧:事件编号(16位)+ 事件值(32位)
    zmq_msg_t msg;
    zmq_msg_init(&msg);
    if (zmq_msg_recv(&msg, monitor, 0) == -1) return -1;
    assert(zmq_msg_more(&msg));

    uint8_t *data = (uint8_t *)zmq_msg_data(&msg);
    uint16_t event = *(uint16_t *)(data);         // 事件编号
    if (value) *value = *(uint32_t *)(data + 2);  // 事件值

    // 第二帧:受影响端点的字符串
    zmq_msg_init(&msg);
    if (zmq_msg_recv(&msg, monitor, 0) == -1) return -1;
    assert(!zmq_msg_more(&msg));

    if (address) {
        uint8_t *data = (uint8_t *)zmq_msg_data(&msg);
        size_t size = zmq_msg_size(&msg);
        *address = (char *)malloc(size + 1);
        memcpy(*address, data, size);
        (*address)[size] = 0;
    }
    return event;
}

事件循环中维护一个 server 计数:握手成功(ZMQ_EVENT_HANDSHAKE_SUCCEEDED)时加一,断开(ZMQ_EVENT_DISCONNECTED)时减一,归零时触发关闭代理:

cpp
// 套接字监控仅适用于 inproc://
int rc = zmq_socket_monitor(_backend, "inproc://monitor", ZMQ_EVENT_ALL);
assert(rc == 0);

zmq::socket_t server_mon(_ctx, ZMQ_PAIR);
std::thread th1([&]() {
    server_mon.connect("inproc://monitor");
    int event = 0;
    while (event != ZMQ_EVENT_MONITOR_STOPPED)
    {
        event = get_monitor_event(server_mon, NULL, NULL);
        switch (event)
        {
            case ZMQ_EVENT_HANDSHAKE_SUCCEEDED:
                ++_serverCnt;
                break;
            case ZMQ_EVENT_DISCONNECTED:
                --_serverCnt;
                if (_serverCnt == 0)
                {
                    // 关闭 proxy,注意资源回收
                }
                break;
            default:
                break;
        }
    }
});

一个实现上的提醒:示例代码里关闭 proxy 用的是 pthread_cancel(),这有可能导致 ZeroMQ 上下文崩溃,仅作演示。生产实现建议改用 zmq_proxy_steerable() 代替 zmq_proxy(),向其 control 端点发送 TERMINATE 指令来优雅终止代理。

完整示例 ​

以下代码实现多对多的异步客户端/服务器模式:proxy 统计握手成功的 server 数量,归零后关闭代理。三段代码可分别编译运行。

proxy ​

cpp
#include <czmq.h>

#include <atomic>
#include <condition_variable>
#include <iostream>
#include <thread>
#include <zmq.hpp>

class Proxy
{
public:
    Proxy()
        : _ctx(1),
          _frontend(_ctx, ZMQ_ROUTER),
          _backend(_ctx, ZMQ_DEALER),
          _capture(_ctx, ZMQ_DEALER),
          _eventCollect(_ctx, ZMQ_PAIR),
          _msgCollect(_ctx, ZMQ_DEALER)
    {
    }

    void init()
    {
        try
        {
            // server 归零后异步关闭 proxy
            std::thread t([&]() {
                while (true)
                {
                    std::unique_lock<std::mutex> lck(_mtx);
                    _cond.wait(lck);
                    // 等待 server 重启窗口
                    std::this_thread::sleep_for(std::chrono::seconds(3));
                    if (_serverCnt == 0)
                    {
                        // 生产环境建议用 proxy_steurable() + TERMINATE 指令
                        pthread_cancel(_proxyThread.native_handle());
                        break;
                    }
                }
            });
            t.detach();

            // 套接字监控仅适用于 inproc://
            int rc = zmq_socket_monitor(_backend, "inproc://monitor", ZMQ_EVENT_ALL);
            assert(rc == 0);

            std::thread th1([&]() {
                _eventCollect.connect("inproc://monitor");
                int event = 0;
                while (event != ZMQ_EVENT_MONITOR_STOPPED)
                {
                    event = get_monitor_event(_eventCollect, NULL, NULL);
                    switch (event)
                    {
                        case ZMQ_EVENT_HANDSHAKE_SUCCEEDED:
                            ++_serverCnt;
                            break;
                        case ZMQ_EVENT_DISCONNECTED:
                            --_serverCnt;
                            if (_serverCnt == 0)
                            {
                                _cond.notify_one();
                            }
                            break;
                        default:
                            break;
                    }
                }
            });
            _monitorThread.swap(th1);

            _frontend.setsockopt(ZMQ_SNDHWM, 10000);
            _backend.setsockopt(ZMQ_HEARTBEAT_IVL, 1000);
            _backend.setsockopt(ZMQ_HEARTBEAT_TIMEOUT, 5000);

            _frontend.bind("tcp://*:5570");
            _backend.bind("tcp://*:5571");
            _capture.bind("inproc://capture");

            // 捕获消息的线程:旁路统计经过 proxy 的流量
            std::thread th2([&]() {
                _msgCollect.connect("inproc://capture");
                int cnt = 0;
                while (true)
                {
                    zmq::message_t identity;
                    zmq::message_t msg;
                    _msgCollect.recv(&identity);
                    _msgCollect.recv(&msg);
                    std::cout << "count:" << ++cnt << std::endl;
                }
            });
            _captureThread.swap(th2);
        } catch (std::exception &e)
        {
            std::cerr << e.what();
        }
    }

    void join()
    {
        _proxyThread.join();
        _frontend.close();
        _backend.close();
        _capture.close();

        pthread_cancel(_captureThread.native_handle());
        _captureThread.join();
        _msgCollect.close();
        _monitorThread.join();
        _eventCollect.close();
        _ctx.close();
    }

    void run()
    {
        try
        {
            std::thread th([&]() {
                zmq::proxy(static_cast<void *>(_frontend), static_cast<void *>(_backend),
                           static_cast<void *>(_capture));
            });
            _proxyThread.swap(th);
            std::cout << "proxy running..." << std::endl;
        } catch (std::exception &e)
        {
            std::cerr << e.what();
        }
    }

private:
    int get_monitor_event(void *monitor, int *value, char **address)
    {
        zmq_msg_t msg;
        zmq_msg_init(&msg);
        if (zmq_msg_recv(&msg, monitor, 0) == -1) return -1;
        uint8_t *data = (uint8_t *)zmq_msg_data(&msg);
        uint16_t event = *(uint16_t *)(data);
        if (value) *value = *(uint32_t *)(data + 2);

        zmq_msg_init(&msg);
        if (zmq_msg_recv(&msg, monitor, 0) == -1) return -1;

        if (address)
        {
            uint8_t *data = (uint8_t *)zmq_msg_data(&msg);
            size_t size = zmq_msg_size(&msg);
            *address = (char *)malloc(size + 1);
            memcpy(*address, data, size);
            (*address)[size] = 0;
        }
        return event;
    }

private:
    zmq::context_t _ctx;
    zmq::socket_t _frontend;
    zmq::socket_t _backend;
    zmq::socket_t _capture;
    zmq::socket_t _eventCollect;
    zmq::socket_t _msgCollect;
    std::atomic<int> _serverCnt{0};
    std::thread _proxyThread;
    std::thread _captureThread;
    std::thread _monitorThread;
    mutable std::mutex _mtx;
    std::condition_variable _cond;
};

int main()
{
    Proxy proxy;
    proxy.init();
    proxy.run();
    proxy.join();
}

server ​

cpp
#include <czmq.h>

#include <thread>
#include <zmq.hpp>

class Server
{
public:
    Server() : _ctx(1), _serverSocket(_ctx, ZMQ_DEALER) {}

    void start()
    {
        _serverSocket.connect("tcp://localhost:5571");

        try
        {
            while (true)
            {
                zmq::message_t identity;
                zmq::message_t msg;
                zmq::message_t copied_id;

                _serverSocket.recv(&identity);
                _serverSocket.recv(&msg);

                // 每个请求回复 5 条,验证多对多转发
                for (int reply = 0; reply < 5; ++reply)
                {
                    copied_id.copy(&identity);
                    _serverSocket.send(copied_id, ZMQ_SNDMORE);
                    char response_string[50] = {};
                    sprintf(response_string, "%s, response #%d", msg.to_string().c_str(), reply);
                    _serverSocket.send(response_string, strlen(response_string), ZMQ_DONTWAIT);
                    std::this_thread::sleep_for(std::chrono::milliseconds(1));
                }
            }
        } catch (std::exception &e)
        {
        }
    }

private:
    zmq::context_t _ctx;
    zmq::socket_t _serverSocket;
};

int main()
{
    Server ser;
    ser.start();
}

client ​

cpp
#include <czmq.h>

#include <iostream>
#include <thread>
#include <zmq.hpp>

#define within(num) (int)((float)((num)*random()) / (RAND_MAX + 1.0))

class Client
{
public:
    Client() : _ctx(1), _clientSocket(_ctx, ZMQ_DEALER) {}
    void start()
    {
        char identity[10] = {};
        sprintf(identity, "%04X-%04X", within(0x10000), within(0x10000));
        _clientSocket.setsockopt(ZMQ_IDENTITY, identity, strlen(identity));
        _clientSocket.setsockopt(ZMQ_RCVHWM, 5000);
        _clientSocket.connect("tcp://localhost:5570");

        zmq::pollitem_t items[] = {{_clientSocket, 0, ZMQ_POLLIN, 0}};
        int request_nbr = 0;
        try
        {
            std::thread th([&]() {
                while (true)
                {
                    zmq::poll(items, 1, 1);
                    if (items[0].revents & ZMQ_POLLIN)
                    {
                        s_dump(_clientSocket);
                    }
                }
            });
            int cnt = 100;
            do
            {
                char request_string[16] = {};
                sprintf(request_string, "request #%d", ++request_nbr);
                std::cout << "send " << request_string << std::endl;
                _clientSocket.send(request_string, strlen(request_string), ZMQ_DONTWAIT);
                std::this_thread::sleep_for(std::chrono::seconds(1));
            } while (--cnt);

            th.join();
        } catch (std::exception &e)
        {
        }
    }

private:
    zmq::context_t _ctx;
    zmq::socket_t _clientSocket;
};

int main()
{
    Client cli;
    cli.start();
}

要点 ​

  • server 进程退出由内核完成 TCP 四次挥手,socket monitor 的 ZMQ_EVENT_DISCONNECTED 就能感知;网络中断必须靠 ZMTP 心跳超时判定,两种场景可以统一收敛到 DISCONNECTED 事件上处理。
  • 心跳三选项中,TIMEOUT 依赖 IVL;取消超时不要求收到 PONG,任何流量都算数;TTL 是告知远端的超时,向下取整到 0.1 秒。
  • 心跳选项须在 bind/connect 之前设置,只影响之后建立的连接。
  • 关闭代理优先使用 zmq_proxy_steerable() + TERMINATE 指令,避免对 ZeroMQ 线程使用 pthread_cancel()。

参考 ​

最近更新

基于 VitePress 构建